Java NIO - Netty 的工作模型

案例一:BIO Echo

在 BioSocketServer 中,每当 serverSocket.accept() 监听到一个新连接,就会将其打包成一个任务丢给线程池。这是 一连接一线程(或线程池异步处理) 的 BIO 模型。线程池的线程依然阻塞在 br.readLine() 上。如果客户端连上后不发数据(长连接保持),该线程就会被无限期挂起。面对海量连接时,即使是多线程池也会瞬间耗尽。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
import java.io.BufferedReader;
import java.io.IOException;
import java.io.InputStreamReader;
import java.io.OutputStream;
import java.net.ServerSocket;
import java.net.Socket;
import java.nio.charset.StandardCharsets;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

public class BioSocketServer {
public static void main(String[] args) throws IOException {
ServerSocket serverSocket = new ServerSocket(8888);
ExecutorService executor = Executors.newFixedThreadPool(Runtime.getRuntime().availableProcessors());

try {
while (true) {
final Socket socket = serverSocket.accept();
executor.execute(() -> socketHandler(socket));
}
} finally {
serverSocket.close();
}
}

private static void socketHandler(Socket socket) {
try (socket;
BufferedReader br = new BufferedReader(new InputStreamReader(socket.getInputStream(), StandardCharsets.UTF_8));
OutputStream out = socket.getOutputStream()) {

String line;
while ((line = br.readLine()) != null) {
System.out.println(line);

out.write("哈哈 收到了你的消息\n".getBytes(StandardCharsets.UTF_8));
out.flush();

if ("quit".equalsIgnoreCase(line.trim())) {
break;
}
}
} catch (IOException e) {
throw new RuntimeException(e);
}
}
}


案例二:NIO 群发聊天(单线程Reactor模型)

问题说明

为了解决阻塞问题,我们在 GroupChatServer 中引入了 Selector。这是 I/O 多路复用模型。一个主线程通过 selector.select() 统一监听所有通道的 ACCEPT 和 READ 事件。但是也存在以下两个主要的问题:

  • 单线程压力过大:listen() 方法中,主线程既处理客户端接入(isAcceptable),又负责数据的读取和遍历转发(isReadable)。如果某个客户端发送的消息很大,或者转发逻辑(sentMsgToOthers)涉及到数据库查询或耗时计算,整个 Selector 轮询就会被卡死,导致其他客户端的连接或消息无法及时响应。
  • 编码极其繁琐:需要小心翼翼地处理 iter.remove() 规避重复触发,需要手动处理 count == -1 规避 CPU 100% 飙升,还要手动处理各种半包、断开连接的异常。


案例代码

Server 端

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
import lombok.extern.slf4j.Slf4j;
import java.io.IOException;
import java.net.InetSocketAddress;
import java.nio.ByteBuffer;
import java.nio.channels.*;
import java.nio.charset.StandardCharsets;
import java.util.Iterator;
import java.util.Set;

@Slf4j
public class GroupChatServer {

private ServerSocketChannel serverSocketChannel;
private static final int PORT = 8888;
private Selector selector;

public GroupChatServer() {
try {
serverSocketChannel = ServerSocketChannel.open();
serverSocketChannel.bind(new InetSocketAddress(PORT));
serverSocketChannel.configureBlocking(false);

selector = Selector.open();
serverSocketChannel.register(selector, SelectionKey.OP_ACCEPT);
log.info("[{}] 群聊服务端已启动,监听端口:{}", Thread.currentThread().getName(), PORT);
} catch (IOException e) {
log.error("新建 GroupChatServer 失败", e);
}
}

public void listen() {
try {
while (true) {
log.info("[{}] 等待 select...", Thread.currentThread().getName());
int selectedCount = selector.select(); // 阻塞式的方法
if (selectedCount > 0) {
Set<SelectionKey> selectedKeys = selector.selectedKeys();
Iterator<SelectionKey> iter = selectedKeys.iterator();
while (iter.hasNext()) {
SelectionKey key = iter.next();
iter.remove(); // 立即删除 key 防止重处理

if (key.isAcceptable()) {
SocketChannel socketChannel = serverSocketChannel.accept();
socketChannel.configureBlocking(false); ////
socketChannel.register(selector, SelectionKey.OP_READ);
log.info("[{}] {} 上线...", Thread.currentThread().getName(), socketChannel.getRemoteAddress());
}

if (key.isReadable()) {
readData(key);
}
}
}
}
} catch (IOException e) {
log.info("服务端监听异常", e);
}
}

private void readData(SelectionKey key) {
SocketChannel channel = (SocketChannel) key.channel();
try {
// 读数据
ByteBuffer buffer = ByteBuffer.allocate(1024);
int count = channel.read(buffer);
if (count > 0) {
// 使用 flip 和指定长度构建字符串,防止乱码和脏数据
buffer.flip();
String msg = new String(buffer.array(), 0, buffer.remaining(), StandardCharsets.UTF_8);
log.info("[{}] From 客户端:{}", Thread.currentThread().getName(), msg);
// 向其他客户端转发数据
sentMsgToOthers(msg, channel);
} else if (count == -1) { // 当客户端正常关闭连接(比如直接关闭进程)时
// 处理客户端正常下线,防止死循环导致 CPU 100%
closeChannel(key, channel);
}
} catch (IOException e) {
// 处理客户端异常掉线
closeChannel(key, channel);
}
}

private void sentMsgToOthers(String msg, SocketChannel self) throws IOException {
Set<SelectionKey> keys = selector.keys(); // 遍历所有注册到 selector 上的 SocketChannel 并且排除自己
for (SelectionKey key : keys) {
// 确保 key 有效再处理
if (!key.isValid()) {
continue;
}
SelectableChannel targetChannel = key.channel();
if (targetChannel instanceof SocketChannel channel && targetChannel != self) {
// 消息包装成 buffer 并发到 channel
ByteBuffer byteBuffer = ByteBuffer.wrap(msg.getBytes(StandardCharsets.UTF_8));
channel.write(byteBuffer);
}
}
}

private void closeChannel(SelectionKey key, SocketChannel channel) {
try {
if (channel != null) {
log.warn("[{}] {} 离线了...", Thread.currentThread().getName(), channel.getRemoteAddress());
key.cancel();
channel.close();
}
} catch (IOException ex) {
log.error("关闭通道异常", ex);
}
}

public static void main(String[] args) {
new GroupChatServer().listen();
}
}

Client 端

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
import lombok.extern.slf4j.Slf4j;
import java.io.IOException;
import java.net.InetSocketAddress;
import java.nio.ByteBuffer;
import java.nio.channels.SelectionKey;
import java.nio.channels.Selector;
import java.nio.channels.SocketChannel;
import java.nio.charset.StandardCharsets;
import java.util.Iterator;
import java.util.Scanner;
import java.util.Set;

@Slf4j
public class GroupChatClient {

private static final String HOST = "127.0.0.1";
private static final int PORT = 8888;
private SocketChannel socketChannel;
private Selector selector;
private String username;

public GroupChatClient() {
try {
selector = Selector.open();
socketChannel = SocketChannel.open(new InetSocketAddress(HOST, PORT));
socketChannel.configureBlocking(false);
socketChannel.register(selector, SelectionKey.OP_READ);
username = socketChannel.getLocalAddress().toString().substring(1);
log.info("[{}] 客户端 {} 已经准备好了...", Thread.currentThread().getName(), username);
} catch (IOException e) {
log.error("新建客户端失败", e);
}
}


public void sendMsg(String msg) {
msg = String.format("%s 说:%s", username, msg);
try {
log.info("[{}] 发送消息:{}", Thread.currentThread().getName(), msg);
socketChannel.write(ByteBuffer.wrap(msg.getBytes(StandardCharsets.UTF_8)));
} catch (IOException e) {
log.error("发送消息失败", e);
}
}

public void readMsg() {
try {
int readChannelCount = selector.select(); // 这个方法是阻塞式的
if (readChannelCount > 0) {
Set<SelectionKey> selectedKeys = selector.selectedKeys();
Iterator<SelectionKey> iter = selectedKeys.iterator();
while (iter.hasNext()) {
SelectionKey key = iter.next();
iter.remove(); // 马上移除
if (key.isReadable()) {
SocketChannel channel = (SocketChannel) key.channel();
ByteBuffer byteBuffer = ByteBuffer.allocate(1024);
int count = channel.read(byteBuffer);
if (count > 0) {
// 准确截取有效字节,杜绝尾部空数据
byteBuffer.flip();
String msg = new String(byteBuffer.array(), 0, byteBuffer.remaining(), StandardCharsets.UTF_8);
log.info("[{}] 收到 {} 消息:{}", Thread.currentThread().getName(), channel.getRemoteAddress(), msg);
}
}
}
}
} catch (IOException e) {
log.error("读取消息异常", e);
}
}

public static void main(String[] args) {
GroupChatClient groupChatClient = new GroupChatClient();

// 启动一个线程每个3秒读取数据
new Thread(() -> {
while (true) {
groupChatClient.readMsg();
}
}).start();

// 发送数据给服务端
Scanner scanner = new Scanner(System.in);
while (scanner.hasNextLine()) {
String s = scanner.nextLine();
groupChatClient.sendMsg(s);
}
}
}


Netty 的工作模型

Netty 的设计

为了完美解决上述两个案例的痛点,Netty 继承并优化了原生 NIO 的设计,形成了标准的主从 Reactor 多线程模型。Netty 模型之所以被公认为高性能网络框架的标杆,正是因为它在架构层面上,把 GroupChatServer 中用单线程硬撑起来的 NIO 轮询,拆解成了“多线程分工协作(主从 Reactor)”与“流水线分层处理(Pipeline)”的工业级形态。这使得开发者能够像写 BIO 一样简单地编写业务代码,同时享受 NIO 带来的超高并发吞吐能力。

  • 将 “接入” 与 “读写” 强制剥离(BossGroup 与 WorkerGroup)。在我们的原生 NIO 中:一个 Selector 既管接入,又管读写。但在 Netty 中:
    • BossGroup 内部的 NioEventLoop 相当于以上代码中的 if (key.isAcceptable()) 分支。它只负责拦截新连接,创建好 Channel 后,立刻将其丢给 WorkerGroup。通常只需要配置 1 个线程。
    • WorkerGroup 内部的 NioEventLoop(默认有多个)相当于以上代码中的 if (key.isReadable()) 分支。它专门负责监听读写事件。默认的线程数是 CPU 核心数 × 2
    • 这样设计的好处是 “前台迎宾” 和 “后厨做菜” 完全独立,互不卡顿。
  • 将 “死循环轮询” 抽象为驱动环(NioEventLoop)。以上两个案例中都写了 while(true) 块。Netty 将这个死循环封装成了 NioEventLoop 的动力环,每次循环固定做三件事(NioEventLoop 这三件事完全在同一个线程之中):
    • select():监听网络事件(等数据)。
    • processSelectedKeys():捕获网络事件,触发底层的 I/O 读写。
    • runAllTasks():执行任务队列中的非 I/O 异步任务(例如延时任务或定时任务)。
  • 将 “数据处理与业务” 高度解耦(ChannelPipeline 机制)。在我们的原生 NIO 中,读取到 Buffer 后,数据转换、群发逻辑全部挤在 readData 方法里。在 Netty 中:引入了 ChannelPipeline(责任链模式)。
    • 当 Worker 线程触发读事件后,数据会像流水线一样流过 Pipeline。
    • 我们可以不再需要在一个方法里写完所有逻辑。写一个 DecoderHandler 专门负责把二进制转换成字符串(自动解决黏包问题);再写一个 ChatHandler 专门负责类似 sentMsgToOthers 的群发工作。


设计的精髓

  • 绝对的无锁化(串行设计): 在 Netty 中,一个 Channel 在创建后,整个生命周期内的所有 I/O 操作和 Pipeline 处理器,都由同一个固定的 NioEventLoop(单线程)来执行。这意味着由于没有多线程并发去抢占同一个通道,该通道内部完全不需要加锁,不仅消除了死锁隐患,更极大减少了线程上下文切换的巨大开销。
  • 非阻塞的 I/O 多路复用: 一个 Worker 线程凭借内部的 Selector,可以同时管理和轮询成千上万个客户端连接。只要没有数据可读,线程就在 select() 上等待,绝不浪费 CPU 资源。
  • 极佳的职责分离: Boss 和 Worker 完全解耦。即使后厨的业务 Handler 正在进行高密度的解码或暂时性的计算,也绝对不会影响前台高并发地接待新用户。


具体的工作模型

NioEventLoop 轮询三部曲

无论在 BossGroup 还是 WorkerGroup 中,最底层的执行驱动都是 NioEventLoop。图中的巨大圆环代表它是一个死循环的单线程任务驱动器。它每次循环都在做三件事(图中底部的 step 1/2/3):

  • step 1: select():调用原生 selector.select(timeout),让线程在操作系统内核层面阻塞,等待感兴趣的网络事件发生。
  • step 2: processSelectedKeys():一旦被唤醒,拿到就绪的 SelectionKey 集合,顺着 Pipeline 链条触发对应的通道事件(如 Accept、Read、Write)。
  • step 3: runAllTasks():处理非纯粹 IO 的自定义异步任务(也就是丢进 taskQueue 里的任务,如我们之前讨论的跨线程写操作或定时心跳任务)。


BossGroup 线程:如何处理 Accept?

BossGroup 通常只需要 1 个线程,即 new NioEventLoopGroup(1),因为它只负责一件事情:拦截连接请求。

  • 启动与监听绑定:当服务端调用 serverBootstrap.bind(8888) 时,Netty 底层会创建一个原生的 ServerSocketChannel,并将其配置为非阻塞模式。接着,将它注册到 BossGroup 唯一的 NioEventLoop 的 Selector 上,监听 SelectionKey.OP_ACCEPT 事件。
  • 拦截与接受连接:当客户端发起 TCP 三次握手时,BossGroup 的线程在 step 1 中被唤醒:
    • 内核响应:进入 step 2,线程拿到就绪的 Key,发现是 OP_ACCEPT。
    • 调用 Accept:底层调用原生 ServerSocketChannel.accept(),拿到专门负责和该客户端后续通信的 SocketChannel。
    • 打包传递:Netty 将这个原生的 SocketChannel 封装成自己的 NioSocketChannel 对象,并以入站事件的形式向后抛出。


WorkerGroup 线程:如何接管并处理 Read/Write?

WorkerGroup 通常包含多个线程(默认是电脑 CPU 核心数 x 2)。它们不负责创建新连接,只负责已有连接的数据读写与业务编解码。

  • 跨组协作:Boss 传球给 Worker,这对应图中那根跨组的绿色箭头(注册 Channel 到 Selector)。
    • 当 BossGroup 拿到 NioSocketChannel 后,会通过轮询算法(Round-Robin)在 WorkerGroup 众多的 NioEventLoop 中挑选一个(假设选到了 NioEventLoop 1)。
    • 由于 Boss 线程和 Worker 线程是隔离的,Boss 线程不能直接操作 Worker 的 Selector。Boss 线程会将“注册新通道” 这一行为打包成一个 Runnable Task,丢进 NioEventLoop 1 的 taskQueue 中。
    • 此时,BossGroup 的接待任务完成,立刻返回 Step 1: select 继续等下一个新客户。
  • Worker 异步消费与监听:
    • NioEventLoop 1 在执行完自己的 IO 轮询后,走到 step 3: runAllTasks(),发现了 Boss 丢过来的注册任务,于是由 Worker 线程执行该任务。
    • 注册读事件:Worker 线程将这个 NioSocketChannel 注册到自己独有的 Selector 上,并注册感兴趣的事件为 SelectionKey.OP_READ。
  • 常态化读写流转:此后,该客户端的所有请求都由 NioEventLoop 1 全权负责。
    • 客户端发来消息,NioEventLoop 1 的 Selector 在 step 1 被唤醒。
    • 进入 step 2,调用原生 SocketChannel.read(ByteBuf) 将报文读入内存。
    • 沿着该通道绑定的 ChannelPipeline 顺流而下,依次流经众多业务 XxxHandler 最终完成相关业务逻辑。


核心线程间的协作与线程安全性

Netty 设计中最精妙的就在于“主从 Reactor 线程之间、以及 Worker 之间是如何完美协作却不需要加锁的”:

  • Boss 与 Worker 的协作:遵循非当前线程提交 Task 的机制。Boss 线程只做 “投递” 动作(往 Worker 的 Concurrent 队列塞任务),不直接操作 Worker 内部组件。Worker 在 step 3 自主消费,实现完全的线程隔离。
  • 同一个 Channel 的读写:无锁串行化(图一机制)。一旦 NioSocketChannel 分配给了 NioEventLoop 1,只要这个连接不断开,它一生所有的网络读、写、协议解析、事件触发,有且仅有 NioEventLoop 1 这一根线程在执行。因此,你的 MyWebSocketHandler 里面不需要加任何多线程并发锁。
  • Worker 多个线程之间:井水不犯河水。每个 NioEventLoop 拥有自己独立的原生 Selector、独立的 taskQueue、独立的线程。不同的客户端被均匀散落在不同的 Worker 线程上,互相之间无任何业务交集,从而彻底避免了传统多线程高并发下的死锁和上下文切换开销。


Pipeline 业务管道

当 Worker 线程在 Step 2 触发底层读取到原始的 ByteBuf 字节数据后,它会将数据送入该通道专属的 ChannelPipeline(责任链) 内部:

  • Decoder(解码器 Handler):负责将原本断断续续的二进制字节流,按照特定协议拼接、拆分成完整的业务对象(解决网络黏包、拆包问题)。
  • Business Logic(业务 Handler):真正的群聊业务代码(例如我们之前写的广播、数据库查询)。
  • Encoder(编码器 Handler):当业务层决定回写数据时,将响应对象重新转为二进制字节流,最终通过 Worker 线程发送给客户端。

Pipeline 的正向和逆向驱动:

  • 正向驱动(Inbound):由 fireChannelRead 领衔,数据从左往右(网络层 $\rightarrow$ 业务层)流动。
  • 逆向驱动(Outbound):由 writeAndFlush 领衔,数据从右往左(业务层 $\rightarrow$ 网络层)流动。
    • ctx.writeAndFlush(msg) —— 从当前节点向前逆推。它是 ChannelHandlerContext 的方法。一旦调用,事件会从当前的 Handler 开始,向前(向左)寻找上一个 OutboundHandler。它会跳过当前 Handler 后续(右边)配置的所有 OutboundHandler。
    • ctx.channel().writeAndFlush(msg) —— 从管道尾部彻底逆推。它是 Channel(或 ctx.pipeline())的方法。一旦调用,事件会直接跳到整个 Pipeline 的最后一站(TailContext),然后从头到尾、完整地向前逆推。它能保证管道中配置的所有 OutboundHandler 都会被百分之百执行。

需要注意的是:在开发自定义协议或编解码时,很多人都会遇到这样的 Bug:“我在 Handler 里调用了 write,控制台不报错,但客户端就是收不到数据,网卡抓包也空空如也。” 这通常是由以下两个关于逆向驱动的误区导致的:

  • 忘记调用 flush():write() 方法的作用只是把对象放进 Netty 的内部发送缓冲区(ChannelOutboundBuffer)中,此时数据根本还没有写到操作系统的 Socket 缓冲区,更没有上自网卡。所以,必须调用 ctx.flush() 才会真正清空缓冲区并发送。或者在驱动时直接使用组合方法 writeAndFlush()。

  • 挂载顺序导致的 “完美擦肩而过”。如果你在服务器启动类中,是这样配置 Pipeline 的:

    1
    2
    3
    // ❌ 错误示范
    pipeline.addLast(new MyWebSocketHandler()); // 1. 先挂业务
    pipeline.addLast(new MyProtocolEncoder()); // 2. 后挂编码器

    此时,当你在 MyWebSocketHandler 中调用 ctx.writeAndFlush(“Hello”) 时,Netty 会从当前位置向前(向左)寻找出站处理器。由于 MyProtocolEncoder 挂在它的后面(右边),逆向驱动事件根本看不见它!数据会绕过编码器,直接以原生的 String 类型一路逆流撞向物理 Socket。由于底层的 Socket 处理器只认识 ByteBuf、不认识 String,数据就会在最后一刻被默默丢弃,导致客户端永远收不到消息。正确做法是永远将 Encoder(编码器)挂在业务 Handler 的前面(左边),或者一律使用 ctx.channel().writeAndFlush() 确保从尾部完整逆推。


TaskQueue

TaskQueue解决了什么问题?

在 Netty 的高性能架构中,每个 NioEventLoop 内部不仅包含了一个多路复用器 Selector,还紧密绑定着一个 taskQueue(任务队列)。这个设计是 Netty 实现异步化、非阻塞、串行化核心机制的灵魂所在。

在传统的 NIO 编程或单线程 Reactor 模型中,I/O 线程(即执行事件循环的线程)是绝对不能被阻塞的。如果在 channelRead 等事件回调中执行了一个耗时较长的操作(例如:高密度计算、复杂的编解码、复杂的数据库查询、或者调用远程的 RPC 服务),整个 I/O 线程就会被卡死在 processSelectedKeys() 这一步。一旦线程卡死,它就无法及时去执行 select() 轮询,导致整个服务无法响应其他客户端的新连接或读写事件。

为了解决“网络 I/O 事件的快速响应”与“耗时任务/异步协同的延迟执行”之间的矛盾,Netty 引入了 taskQueue。它的核心目的是允许开发者将非即时、耗时或者跨线程的任务,安全地提交到 I/O 线程的本地队列中,让 I/O 线程在处理完紧急的网络事件后,以异步非阻塞的方式在当前线程“顺手”消费掉,从而保证主事件循环的极速运转。

NioEventLoop 的每次循环都包含三大固定步骤:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
+--------------------------------------------+
| NioEventLoop 死循环 |
+--------------------------------------------+
|
v
[ Step 1: select() ] <------- 阻塞等网络事件
|
v
[ Step 2: processSelectedKeys() ] <------- 快速处理 I/O 读写
|
v
[ Step 3: runAllTasks() ] <------- 消费 taskQueue 里的任务
|
+-----------------------+

详细运转流程:

  • 任务暂存:当有异步任务被提交给 NioEventLoop 时(无论是当前线程还是外部线程提交),该任务都会被压入 taskQueue 队列。
  • 处理时机:I/O 线程在 Step 2 快速把网络数据的读写处理完毕后,就会立刻来到 Step 3 的 runAllTasks()。
  • 时间控制:为了防止队列里的任务太多,导致 runAllTasks() 执行时间过长,从而耽误了下一轮的 select() 轮询,Netty 引入了 ioRatio(I/O 比例)参数。默认情况下,Netty 会限制执行异步任务的时间,使其与处理 I/O 事件的时间保持合理的平衡,一旦超时,剩余任务会留到下一次事件循环再执行。


TaskQueue 三种典型场景

三种典型的使用场景:

  • 用户程序自定义的普通任务:如果你明知某段业务逻辑非常耗时,你可以在 ChannelHandler 内部将该任务包装好提交到当前通道所绑定的 EventLoop 的 taskQueue 中。

    1
    2
    3
    4
    5
    ctx.channel().eventLoop().execute(() -> {
    // 这是一个非常耗时的业务逻辑
    // 它会被放入 taskQueue,由 I/O 线程在 runAllTasks() 阶段异步消费
    doHeavyBudget();
    });
  • 用户自定义定时任务:有些任务不需要立刻执行,而是需要延迟执行或周期性执行(例如:客户端心跳检测未响应踢出、长连接超时断开、定时刷新本地缓存等)。这类任务在 Netty 底层会先进入一个基于小顶堆实现的 scheduledTaskQueue(定时任务队列)。当定时时间到了,它会被挪入普通的 taskQueue 中,在 runAllTasks() 阶段被执行。

    1
    2
    3
    ctx.channel().eventLoop().schedule(() -> {
    System.out.println("3秒后检测客户端心跳...");
    }, 3, TimeUnit.SECONDS);
  • 非当前 Reactor 线程调用 Channel 的各种方法:这在群聊系统或推送系统中极为常见。例如在推送系统的业务线程池/消息队列消费者线程里,你收到了某条消息,需要向特定的用户推送。你通过用户 ID 找到了对应的 Channel 引用,并调用 channel.write(msg)。

    • 由于 Netty 采用绝对串行化设计,绝对不允许外部线程直接去操作 Channel 底层的 Socket 读写。
    • 当你从外部的业务线程调用 channel.write() 时,Netty 底层会自动判断:调用我的线程不是当前 Channel 绑定的 I/O 线程! 于是,Netty 会默默地把这个 write 操作包装成一个 Task,自动提交到该 Channel 专属的 NioEventLoop 的 taskQueue 中。
    • 最终,该写出操作依然会由对应的 I/O 线程在它的事件循环中安全、异步地消费掉。
  • 注意:在 Netty 的语境下,上面的 “I/O 线程” 和 “NioEventLoop 线程” 指的就是同一个线程。

    1
    2
    3
    4
    5
    6
    7
    8
    9
    10
    11
    // 这就是一个唯一的 NioEventLoop 线程(它就是传说中的 I/O 线程)在疯狂运转:
    while (true) {
    // 1. 监听网络事件(I/O 职责)
    select();

    // 2. 捕获事件,读取或写入底层字节流(I/O 职责)
    processSelectedKeys();

    // 3. 执行任务队列里的异步非 I/O 任务(非 I/O 职责,但由它顺手捎带做完)
    runAllTasks();
    }


TaskQueue 任务很重怎么办

因为 select()、processSelectedKeys() 和 runAllTasks() 这三步是在同一个 I/O 线程(即 NioEventLoop 线程)中严格串行执行的。如果 runAllTasks() 里的任务太重(比如耗时太长),这个线程就会被一直卡在第三步。在它把重任务处理完之前,主循环绝对没有办法回到第一步去执行 select()。

这会导致非常严重的后果:后续所有新客户端的连接(Accept)和已有客户端的数据读写(Read/Write),都会在网络缓冲区中排队等待,无法得到及时响应,整个服务的吞吐量和实时性会瞬间暴跌。

针对这个极其重要的问题,Netty 在底层设计了两套防御机制,同时对开发者也提出了核心开发规范:

  • 机制一:Netty 底层的自我防卫 —— ioRatio。 它是为了防止 taskQueue 里的任务太多、太重把 I/O 线程饿死在第三步而引入的时间调节参数,默认值为 50。

    • 当线程进入第二步 processSelectedKeys() 处理网络读写时,Netty 会启动计时,记录下处理 I/O 事件一共花了多少时间(假设本轮 I/O 花了 10 毫秒)。
    • 进入第三步 runAllTasks() 后,Netty 会根据 ioRatio = 50 计算出允许执行异步任务的时间:$\text{允许执行任务的时间} = \text{I/O 时间} \times \left(\frac{100 - ioRatio}{ioRatio}\right)$,由于默认是 50,计算结果也是 10 毫秒。
    • 接下来,I/O 线程每执行完 taskQueue 里的一个任务,就会检查一下当前时间。一旦发现执行异步任务的总时间已经快满 10 毫秒了,它就会果断强行中断退出 runAllTasks()。
    • 哪怕队列里还有 1000 个任务没执行,它也不管了,立刻进入下一轮死循环,先去执行第一步 select() 看看网络上有没有急需处理的读写事件。被中断的任务则留在队列中,等下一次循环再继续做。
  • 机制二:无法防范的 “单次重任务” 陷阱。虽然有 ioRatio 的保护,但它只能防范“任务数量太多”的情况,却无法防范 “单个任务太重”。因为 Netty 只能在两个任务之间去检查时间是否超时。如果用户提交了一个极其恶劣的任务,比如:

    1
    2
    3
    4
    ctx.channel().eventLoop().execute(() -> {
    // 这是一个单次就要执行 5 秒钟的密集计算或数据库死锁查询
    Thread.sleep(5000);
    });

    当 I/O 线程取出这一个任务并开始执行时,Netty 没办法在任务执行到一半时强行掐断它。这就导致 I/O 线程必须死死地被卡在这个任务上 5 秒钟,这 5 秒内,下一轮的 select() 彻底瘫痪。

基于上述原因,我们在编写 Netty 的 ChannelHandler 业务代码时,必须死死遵守一条铁律:绝对不能在 Netty 的 I/O 线程(NioEventLoop)中做任何耗时的、阻塞的操作!如果你遇到了重任务(如复杂的 DB 查询、第三方 RPC 调用、高密度文件读写),工业上标准的解决方案有下面几种:

  • 方案 A:提交到自定义的外部业务线程池(最常用)。在 channelRead 收到数据后,立刻把数据打包,丢给自己创建的 ThreadPoolExecutor 线程池去处理,让 Netty 的 I/O 线程秒回,继续去执行下一轮循环。

  • 方案 B:在 addLast 时指定 EventExecutorGroup。Netty 允许你在向 Pipeline 添加 Handler 时,额外传入一个专门的非 I/O 线程池。这样,这个 Handler 里的所有回调方法(如 channelRead)都不会由 NioEventLoop 线程执行,而是由这个指定的线程池来跑。

    1
    2
    3
    4
    5
    // 创建一个专用的业务线程池
    EventExecutorGroup businessGroup = new DefaultEventExecutorGroup(16);

    // 在初始化管道时指定该线程池
    pipeline.addLast(businessGroup, new MyHeavyBusinessHandler());
  • 方案C:直接提交任务到 EventExecutorGroup 线程池中。DefaultEventExecutorGroup 不仅能在 addLast 时作为一个整体挂载给 Handler,它本身也继承了 Java 标准的 ScheduledExecutorService 和 Netty 的 EventExecutorGroup 接口。这意味着它自己就是一个纯正的、可以直接通过 execute()、submit() 或 schedule() 提交异步任务的线程池。所以,如果你不想把整个 Handler 的生命周期都交出去,只想在特定的某几行核心代码上“减负”,你完全可以像下面这样写:

    1
    2
    3
    4
    5
    6
    7
    8
    9
    10
    11
    12
    13
    14
    15
    16
    17
    18
    19
    20
    21
    22
    23
    24
    25
    26
    27
    28
    // 在类成员变量或单例中初始化(注意不要在 initChannel 里每次都 new,那会造成线程暴涨!)
    static final EventExecutorGroup businessGroup = new DefaultEventExecutorGroup(16);

    @Override
    protected void channelRead0(ChannelHandlerContext ctx, TextWebSocketFrame msg) {
    // 1. 在 I/O 线程中做一些轻量级的、快速的前置检查(如协议判断)
    String rawText = msg.text();

    // 2. 直接将耗时的“重任务”作为 task 提交给 DefaultEventExecutorGroup 异步执行
    // 注意:跨线程传递了包装 ByteBuf 的 msg,一定要进行 retain/release 处理!
    msg.retain();

    businessGroup.execute(() -> {
    try {
    // 这里已经是 businessGroup 的业务线程在执行了
    System.out.println("当前线程:" + Thread.currentThread().getName()); // 打印出来会是 defaultEventExecutor-X-X

    // 模拟耗时的 DB 慢查询或第三方 RPC
    String result = queryFromDatabase(rawText);

    // 3. 执行完毕后,利用 ctx 写回客户端(Netty 会自动通过 taskQueue 帮我们切回 I/O 线程写出,保证线程安全!)
    ctx.writeAndFlush(new TextWebSocketFrame(result));
    } finally {
    // 4. 业务线程消费完毕,手动释放
    msg.release();
    }
    });
    }


业务链中的用户自定义事件

在 Netty 中,所有的事件分为两类:

  • Outbound(出站事件):由你的代码主动发起的操作,比如 connect()、bind()、write()。它们在 Pipeline 中是从后往前(Tail -> Head)反向传播的。
  • Inbound(入站事件):由底层操作系统或 Netty 内部触发的被动事件,比如数据可读(channelRead)、连接激活(channelActive)、以及用户自定义事件(userEventTriggered)。它们在 Pipeline 中是从前往后(Head -> Tail)正向传播的。

用户自定义事件 userEventTriggered 属于 Inbound 事件,必须 “自上而下” 传播。它只能捕获在它前面的 Handler 或者是 Netty 头部(HeadContext) 抛出的事件。

在实际的开发中,谁实现了 userEventTriggered,谁就必须负责 “往下传”!这是最容易踩坑的地方。Netty 的事件流是靠程序员手动 “交接棒” 来延续的。当你继承 SimpleChannelInboundHandler 或 ChannelInboundHandlerAdapter 并重写了 userEventTriggered 方法时,你就成了这个事件的终点站。默认情况下,如果你不作为,事件到了你这里就彻底终止(被拦截)了,后面的 Handler 将再也收不到这个事件。

❌ 错误示范:导致后面的 Handler 变成 “聋子”。

1
2
3
4
5
6
7
8
@Override
public void userEventTriggered(ChannelHandlerContext ctx, Object evt) throws Exception {
if (evt instanceof IdleStateEvent) {
System.out.println("我处理了超时事件");
// 结束了,事件死在了这里
}
// 哪怕不是超时事件(比如其他的自定义事件),也死在了这里!
}

✅ 正确示范:不管你处理没处理,都要显式 “向下传递”。
为了保证链条不被打断,除非你确定这个事件不需要其他人管了,否则必须在结尾调用事件传递方法。

1
2
3
4
5
6
7
8
9
10
@Override
public void userEventTriggered(ChannelHandlerContext ctx, Object evt) throws Exception {
if (evt instanceof IdleStateEvent) {
System.out.println("我处理了超时事件");
}

// 关键:把棒子传给 Pipeline 中的下一个 Handler
super.userEventTriggered(ctx, evt);
// 或者执行:ctx.fireUserEventTriggered(evt);
}

虽然捕获受方向限制,但是任何一个 Handler 都可以作为源头,发射一个自定义事件。你可以在任何 Handler 的任何地方调用:

1
2
3
4
5
// 从当前 Handler 开始,向后面的 Handler 传播这个事件
ctx.fireUserEventTriggered(new MyCustomEvent());

// 或者从整个 Pipeline 的最头部(Head)开始,向所有 Handler 传播这个事件
ctx.pipeline().fireUserEventTriggered(new MyCustomEvent());